Skip to content

[SPARK-60112][CORE] Track memory consumers by identity in TaskMemoryManager - #59331

Open
viirya wants to merge 2 commits into
apache:masterfrom
viirya:SPARK-60112
Open

viirya wants to merge 2 commits into
apache:masterfrom
viirya:SPARK-60112

Conversation

@viirya

@viirya viirya commented Oct 10, 2026

Copy link
Copy Markdown
Member

What changes were proposed in this pull request?

TaskMemoryManager tracks spillable consumers in a HashSet<MemoryConsumer>, which deduplicates by equals/hashCode. This PR switches it to an identity-based set (Collections.newSetFromMap(new IdentityHashMap<>())), so every consumer instance is tracked regardless of how it defines equality.

The set is only used for add, iteration and clear (no contains/remove), and the spill heuristic already compares the requesting consumer by reference (c == requestingConsumer), so the switch does not affect any other logic. This is the only hash-based collection keyed by MemoryConsumer in the codebase.

Why are the changes needed?

With value-based equality, an equal but distinct consumer is silently dropped from the set. That instance:

  • is never offered for spilling, so a memory request can fail to free memory that is actually spillable; and
  • is missing from the per-consumer memory breakdown in UNABLE_TO_ACQUIRE_MEMORY errors, where its memory shows up as "not attributed to a specific consumer".

This affects master today: HybridRowQueue is a Scala case class extending MemoryConsumer, so two queues created with the same arguments in the same task are deduplicated. #59282 fixes this for HybridRowQueue specifically; as suggested by @dongjoon-hyun in that review, this PR fixes it in one place for all current and future consumers. The two PRs are independent and can be merged in either order.

Does this PR introduce any user-facing change?

Yes, as a bug fix: distinct memory consumers that compare equal (currently HybridRowQueue instances with identical arguments in the same task) are now all eligible for spilling and are all listed in the OOM memory breakdown. There is no API or configuration change.

How was this patch tested?

Added two tests to TaskMemoryManagerSuite using a test consumer with value-based equals/hashCode:

  • equalButDistinctConsumersAreAllSpilled: before this change, the second equal consumer was not spilled (50 B left).
  • equalButDistinctConsumersAreAllInMemoryConsumptionBreakdown: before this change, the second consumer's memory was reported as not attributed to a specific consumer.

Both fail without the fix and pass with it. Also ran the org.apache.spark.memory suites, UnsafeExternalSorterSuite, ShuffleExternalSorterSuite, BytesToBytesMap{On,Off}HeapSuite, ExternalSorterSpillSuite and ExternalAppendOnlyMapSuite, plus core/checkstyle and core/Test/checkstyle.

Was this patch authored or co-authored using generative AI tooling?

Generated-by: Claude Code (Claude Opus 5.5)

This pull request and its description were written by Isaac.

…anager

TaskMemoryManager tracked spillable consumers in a HashSet, so a consumer
with value-based equals/hashCode (for example a Scala case class such as
HybridRowQueue) was deduplicated against an equal but distinct consumer.
The dropped instance was never offered for spilling, and its memory was
reported as not attributed to any consumer in the OOM breakdown.

Track consumers in an identity-based set instead, so every consumer
instance is considered regardless of how it defines equality.

Co-authored-by: Claude Code <noreply@anthropic.com>

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you for making this a general fix, @viirya. Switching consumers to an identity set looks correct to me. consumers is only used for add, iteration and clear, and the code that compares consumers already uses reference equality: c == requestingConsumer in both spill loops, request.consumer == consumer in allocatePage, and trigger == this in HybridQueue.spill. I also traced both new tests against the old HashSet, and both fail there as described. I left 6 inline comments, mostly about test coverage. In summary:

  1. TaskMemoryManagerSuite.java L889: The breakdown test relies on distinct toString names

    Equal consumers whose toString is also value-based (HybridRowQueue on the current master) are still listed as indistinguishable lines.

  2. TaskMemoryManager.java L121: Equal consumers are now retained until the task ends

    Nothing removes an entry from consumers before cleanUpAllAllocatedMemory(), so every equal instance stays reachable and is scanned on every spill.

  3. TaskMemoryManagerSuite.java L908: Two affected paths are not tested

    The page allocation recovery path (spillConsumersForPageAllocation) and self-spill of a consumer that used to be deduplicated.

  4. TaskMemoryManagerSuite.java L869: nit: Placement of the helper class

    The other helper classes are grouped at the top of the suite.

  5. TaskMemoryManager.java L116: The identity contract is only documented on a private field

    Stating it on MemoryConsumer would keep a future hash collection keyed by consumers from bringing the bug back.

  6. TaskMemoryManager.java L186: nit: IdentityHashMap allocates its table eagerly

    The no-arg constructor allocates a 64-slot table for every task.

}

@Override
public String toString() {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ValueEqualConsumer returns a distinct name from toString, so this test passes even though the breakdown cannot tell equal consumers apart when their toString is also value-based. That is the case for HybridRowQueue on the current master: two equal queues are both listed with the same case-class toString, e.g. HybridRowQueue(org.apache.spark.memory.TaskMemoryManager@...,<dir>,1,...,false): 30.0 B. #59282 adds an identity suffix to HybridRowQueue.toString, but only for that class.

Could we make toString identical for both instances here and assert that both amounts are listed (: 20.0 B and : 30.0 B)? Alternatively, renderConsumerBreakdown and logMemoryUsage could append System.identityHashCode to each consumer, which keeps every line distinguishable regardless of how a consumer implements toString.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good point. ValueEqualConsumer now prints the same for both instances, and the test asserts that both ValueEqualConsumer: 20.0 B and ValueEqualConsumer: 30.0 B are listed. I'd rather not change the breakdown format for every consumer in this PR. Most consumers already print a distinct Object.toString, and #59282 covers HybridRowQueue. If we want an identity suffix in general, I can open a separate JIRA.

*/
@GuardedBy("this")
private final HashSet<MemoryConsumer> consumers;
private final Set<MemoryConsumer> consumers;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nothing removes an entry from consumers before cleanUpAllAllocatedMemory(). So every equal-but-distinct instance that the old HashSet collapsed into one now stays reachable until the task ends, together with everything it references, and both spill loops iterate over it on every spill.

This already holds for any other distinct consumer, so it is not a new kind of problem, and a task usually has only a few consumers. Still, it may be worth a separate JIRA to drop a consumer from the set once its usage reaches zero (or when it is closed), instead of only ever adding to it.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed. This is the same lifetime that any other distinct consumer already has, so I'll keep this PR focused. Filed SPARK-60134 to drop a consumer from the set once it no longer holds memory.

c2.use(50);

// Both equal consumers must be offered for spilling to satisfy the request.
c3.use(100);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

These tests cover another consumer spilling through acquireExecutionMemory. Two other paths also missed the deduplicated consumer before this change, and they are not covered:

  • The page allocation recovery path. spillConsumersForPageAllocation builds its own candidate map from consumers.
  • Self-spill, where the requester is the equal-but-distinct consumer. It was never in the set, so it never got the key 0 and was never asked to spill itself. For example, c1 acquires first and then frees its memory, c2.use(100) fills the limit, and then c2.use(50) gets nothing without this fix.

A small case for each would keep a later change to either loop from bringing the bug back.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added equalButDistinctConsumerCanSelfSpill (your scenario) and equalButDistinctConsumersAreAllSpilledForPageAllocation. In the second test the grant succeeds without spilling, so only spillConsumersForPageAllocation spills. Both fail with the old HashSet.

* A consumer with value-based equality: two instances on the same manager and memory mode are
* equal (and share a hash code) even though they track memory independently.
*/
private static final class ValueEqualConsumer extends TestMemoryConsumer {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: The other helper classes in this suite (TestAllocator through NonSpillingAllocatingConsumer) are grouped at the top of the class. Could we move ValueEqualConsumer next to them?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done, moved next to the other helpers.


/**
* Tracks spillable memory consumers.
* Tracks spillable memory consumers. Consumers are tracked by identity rather than by

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This Javadoc is the only place that says consumers are tracked by identity, and it is on a private field. Since the goal is to cover all current and future consumers, it may be worth stating the same in the MemoryConsumer class Javadoc as well, e.g. that consumers are tracked by identity and equals/hashCode must not be relied on for that. Then someone who later adds a hash-based collection keyed by MemoryConsumer (for example, for per-consumer statistics) would see the requirement where they look.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done, added a paragraph to the MemoryConsumer class Javadoc.

this.tungstenMemoryAllocator = tungstenMemoryAllocator;
this.taskAttemptId = taskAttemptId;
this.consumers = new HashSet<>();
this.consumers = Collections.newSetFromMap(new IdentityHashMap<>());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: The no-arg IdentityHashMap constructor eagerly allocates a 64-slot table (for the default expected maximum size of 21) for every TaskMemoryManager, i.e. once per task, while the previous HashSet allocated its table lazily on the first add. Since a task rarely has more than a few consumers, a smaller expected size such as new IdentityHashMap<>(4) would avoid most of that. It is minor either way.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done, it now uses new IdentityHashMap<>(4).

- Make ValueEqualConsumer print the same for equal instances and assert both
  amounts are listed in the breakdown; move it next to the other helpers.
- Add tests for self-spill and the page allocation recovery path.
- Document identity tracking in the MemoryConsumer Javadoc.
- Size the IdentityHashMap for a few consumers instead of the default.

Co-authored-by: Claude Code <noreply@anthropic.com>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants